perf(channels): shard outbound dispatch per conversation - #1464
Conversation
A single dispatch goroutine delivered every reply in the process — every channel, every user — each Send blocking the next. One media upload of a few seconds froze all outbound traffic for that long, and sustained throughput was capped at one Send round-trip regardless of lane capacity. The ordering that loop actually protected is per-conversation: a run emits block replies, retry notices and a final answer to one ChatID, and those must arrive in that order. A global order was never required. Dispatch now shards on channel+ChatID. A conversation always maps to the same shard and is delivered serially there, so the ordering guarantee is unchanged, while unrelated conversations proceed in parallel. GOCLAW_OUTBOUND_SHARDS tunes the worker count (default 8). Temp media needed a real fix to go with it. Serial dispatch deduped it implicitly — the first delivery removed the file, and a later message carrying the same path found os.Stat failing and skipped it. Across concurrent shards that check alone is a TOCTOU race, so claims are now taken atomically and released after the send. Also makes the pool limits that bound this path configurable, defaults unchanged: GOCLAW_PG_MAX_OPEN_CONNS / GOCLAW_PG_MAX_IDLE_CONNS, and GOCLAW_HTTP_MAX_IDLE_CONNS / GOCLAW_HTTP_MAX_IDLE_CONNS_PER_HOST. The per-host limit binds when one provider takes nearly all traffic: every request past it pays a fresh TCP+TLS handshake, which a deployment raising LaneMain needs to be able to lift.
clark-cant
left a comment
There was a problem hiding this comment.
Maintainer Review — github-maintain
Summary: Shards outbound dispatch by channel+ChatID to eliminate the single-goroutine bottleneck while preserving per-conversation ordering. Also fixes a TOCTOU race in temp media deduplication and makes connection pool limits configurable.
Risk level: Low-Medium — touches core dispatch path but well-tested with comprehensive regression coverage.
Mandatory gates
| Gate | Result |
|---|---|
| Duplicate / prior work | Clear — no prior PR addresses dispatch sharding |
| Project standards | Excellent — follows existing patterns, well-documented |
| Strategic necessity | High — removes a real throughput ceiling that blocks all outbound traffic behind the slowest send |
Code review
Correctness: Solid. The shard key (channel+ChatID) correctly preserves the ordering guarantee that runs depend on — block replies, retry notices, and final answer to one ChatID must arrive in order. Unrelated conversations proceed in parallel as intended.
Temp media fix: The atomic claim mechanism (sync.Map LoadOrStore) properly closes the TOCTOU race that serial dispatch masked implicitly. The claim set is bounded by the send duration, so it won't grow unbounded.
Pool configurability: Sensible defaults unchanged. The documentation explaining when each limit binds (per-host for single-provider deployments, per-query not per-run for Postgres) is excellent operator guidance.
Test coverage: Comprehensive — 50-message ordering test, slow-chat-does-not-block regression test (the exact scenario this change exists for), 16-goroutine temp media race test, non-temp passthrough verification.
Quality:
- Code is well-structured with clear separation (shard index calculation, claim/release helpers)
- Comments explain why not just what
- No anti-slop patterns: no unnecessary abstractions, no defensive paranoia
- Documentation in
docs/00-architecture-overview.mdis operator-facing and actionable
Verdict: APPROVE ✅
This is a high-quality performance improvement that addresses a real bottleneck. The code is well-tested, well-documented, and follows project conventions. Safe to merge once CI is green.
Reviewed by github-maintain at 2026-07-31T14:58Z
Problem
Manager.dispatchOutboundran as a single goroutine callingchannel.Send()synchronously. Every reply in the process — every channel, every user — queued behind the slowest send.Two consequences:
LaneMaincapacity was configured.Why serialization was not required
The ordering that loop actually protected is per-conversation: a run emits block replies, a possible
Provider busy, retrying...notice, and the final answer to oneChatID, and those must arrive in that order. There are 41PublishOutboundcall sites, several of which target the same chat within one run.A global order was never needed. Ordering across different conversations carries no meaning.
Change
Dispatch now shards on
channel + ChatID:GOCLAW_OUTBOUND_SHARDStunes worker count (default 8).The channel name is part of the shard key because
ChatIDs are only unique within a channel.Temp media needed a real fix alongside it
Serial dispatch deduped temp media implicitly: the first delivery removed the file, and a later message carrying the same path found
os.Statfailing and skipped it. The existing comment ("already sent by another dispatch") states the assumption.Across concurrent shards that check alone is a TOCTOU race — two shards can both stat a live file and upload it twice. Claims are now taken atomically and released after the send. Non-temp (workspace) media passes through unclaimed, since it is never removed after delivery.
Also included
Pool limits that bound this same path, made configurable with defaults unchanged:
GOCLAW_PG_MAX_OPEN_CONNSGOCLAW_PG_MAX_IDLE_CONNSGOCLAW_HTTP_MAX_IDLE_CONNSGOCLAW_HTTP_MAX_IDLE_CONNS_PER_HOSTThe per-host limit binds when one provider takes nearly all traffic: every concurrent request past it pays a fresh TCP+TLS handshake instead of reusing a pooled connection. A deployment raising
LaneMainneeds to be able to lift it. Self-hosted / in-network providers have no reason to be throttled at 10.Tests
internal/channels/dispatch_shard_test.go:Sendis held hostage and the other must still be deliveredVerification
go build ./...,go build -tags sqliteonly ./...,go vet ./...— cleango test -racegreen forchannels,scheduler,store/pg,providersNote:
internal/channels/telegramTestDisplayWidthfails on this branch, but it fails identically on a pristineorigin/devworktree — pre-existing and unrelated.